[REC-AI-022] Реализовать семантический DLQ-инкубатор для неструктурированных ответов LLM
Обработка Semantic Garbage ошибок ИИ-моделей Ollama без блокировки конвейера Kafka
- Репозиторий / Компонент:
receipt-processor-service/package worker. - Тип задачи: Инфраструктурная / Обработка ошибок ИИ (AI Data Pipeline).
- Связанные документы:
receipt-processor.qmd(Раздел 3.1 Инвариант изоляции ИИ-сбоев). - Статус: Готово к реализации
Описание задачи:
Текущая архитектура
OCRWorkerиFilterWorkerкорректно обрабатывает инфраструктурные сбои (сетевая недоступность Ollama). Однако существует критическая уязвимость типа Semantic Garbage (Семантический мусор): если сервер ИИ доступен, но локальная модельLlama VisionилиQwenсгенерировала битый, невалидный JSON, методjson.Unmarshalвыбрасывает ошибку парсинга. В текущей реализации воркер прерывает выполнение и оставляет сообщение в Кафке, что вызывает бесконечный цикл повторов (infinite retry Loop) и намертво блокирует чтение всей партиции Kafka для чеков других пользователей.Необходимо внедрить семантический паттерн Dead Letter Queue (DLQ) с автоматическим транзитом бракованных ИИ-структур в очередь модерации.
Инструкция по шагам:
- Разделение типов ошибок: Внутри методов
ProcessMessageворкеров разделить обработку ошибок на два класса:- Инфраструктурные (ошибки сети
http.Do): возвращатьerror, сохраняя сообщение в топике. - Семантические (ошибки парсинга
json.Unmarshalили пустой JSON): перехватывать и отправлять в изолятор.
- Инфраструктурные (ошибки сети
- Реализация транзита в DLQ: При фиксации семантического сбоя, воркер должен извлечь исходное сообщение, обернуть его в метаструктуру с указанием типа сбоя (
ocr_structural_broken_json) и опубликовать в выделенный топик Кафкиreceipts-dead-letter-queue. После публикации воркер обязан выполнитьCommit Offset, освобождая партицию для следующих сообщений. - Автоматическая интеграция с очередью модерации: Создать выделенный консьюмер
DLQModerationConsumer, который слушает топик отстоя, извлекает сырой текст сбойного ответа ИИ и атомарно записывает его в таблицуvoice_moderation_queue(или смежнуюreceipt_moderation_queue) со статусомpendingи детальным логом для ручного исправления оператором.
- Разделение типов ошибок: Внутри методов